Repository navigation
Databricks: detect Delta table version changes (#74195). - #74225
goyaladitay11 wants to merge 2 commits into
Conversation
|
Hey @zozo123, Addressed Phase 1 of #74195 with a fully deferrable DatabricksDeltaTableVersionSensor! Highlights: Full support for UnityTableIdentity + table_name |
|
Thanks @goyaladitay11 — really appreciate you jumping on Phase 1 of #74195. The shape is right: a deferrable A few notes before this is mergeable, in the spirit of keeping Unity assets trustworthy for Airflow + Databricks SQL / jobs users:
Happy to re-review once this sits on top of #74191’s identity types. Thanks again for taking this on. |
0b77124 to
9bb3d0e
Compare
|
Thanks for the clear feedback @zozo123! I have addressed all three points:
All 46 unit tests for sensor, trigger, and assets are passing. Ready for your re-review whenever you have a chance! |
zozo123
left a comment
There was a problem hiding this comment.
Thanks for the follow up. Observation/coalescing semantics and Phase 1 docs look good. Re-reviewed at 9bb3d0e, four things before approval:
- Workspace check (P1). The sensor observes via
databricks_conn_idbut stampsunity_table.hoston provenance and the outlet (L174, L249). Please reuse the COPY INTO guard (here):(hook.host or "").lower() != self.unity_table.host. - Deferral drops config (P2).
session_configuration,http_headers,client_parameters,hook_paramsand query tags reach the sync hook but not the trigger (defer, trigger). Please serialize them into the trigger and use them in_get_hook. - Case folding (P2).
UnityTableIdentitylowercases, but L152 compares exactly, soMain.Default.Customersfails. Lowercase the parts first, like COPY INTO. - Pinned read docs (P2). The example reads default XCom as a dict; sync returns
None, deferred returns an int. Useti.xcom_pull(task_ids="sensor", key="delta_table_version")["version"]and fix thereturn_valuenote on L88.
Also please rebase onto main so the #74191 commits drop out, and add tests for host mismatch and mixed case table_name. Happy to re-review after.
9bb3d0e to
49f321b
Compare
|
Thanks @zozo123! All four points have been addressed, verified locally (all 1,043 unit tests passing), and pushed: 1.Workspace check (P1): Added the COPY INTO workspace host guard Ready for your re-review! |
zozo123
left a comment
There was a problem hiding this comment.
Thanks, re-reviewed at 49f321b. The four earlier points and the rebase look good. One remaining P2 before approval:
Deferral swallows query errors. Sync _get_latest_version raises AirflowException on TABLE_OR_VIEW_NOT_FOUND or permission failures. The trigger except in run only logs a warning and retries until timeout, so a bad table name can burn the full sensor timeout. Please yield an error TriggerEvent for permanent query failures (missing table, access denied) and keep retry only for transient errors, matching the common SQL trigger. Related: restore baseline_version from the trigger event in execute_complete before _push_provenance so deferred enrollment provenance stays accurate after resume.
Happy to approve right after that.
49f321b to
4c4461c
Compare
|
Thanks @zozo123! Both items have been addressed: Permanent Query Error Handling in Trigger: Added _is_permanent_error() to detect permanent query failures (TABLE_OR_VIEW_NOT_FOUND, PERMISSION_DENIED, ACCESS_DENIED, etc.). The trigger now immediately yields an error TriggerEvent and exits on permanent failures instead of retrying until timeout, preserving retries only for transient errors. Added unit tests for missing table, permission denied, and transient error retry. |
zozo123
left a comment
There was a problem hiding this comment.
Re-reviewed at 4c4461c. Both remaining P2s look good:
- Trigger
_is_permanent_errornow yields an errorTriggerEventand returns on missing-table / permission failures; only transient errors are retried (with tests for both paths). - Success events carry
baseline_version, andexecute_completerestores it before_push_provenance(covers deferred enrollment).
Earlier workspace host check, deferral config pass-through, case-folding, and pinned-read docs/delta_table_version XCom key still look intact. LGTM for Phase 1.
4c4461c to
bc2fb2e
Compare
|
Thank you so much for the thorough and thoughtful reviews throughout this PR, @zozo123! Really enjoyed working on this. |
bc2fb2e to
1d75994
Compare
| DatabricksDeltaTableVersionSensor (Phase 1) | ||
| =========================================== | ||
|
|
||
| Use the :class:`~airflow.providers.databricks.sensors.databricks_delta_table.DatabricksDeltaTableVersionSensor` | ||
| to monitor a Delta table in Databricks and detect new commit versions. When a newer version is observed, | ||
| the sensor succeeds, pushes provenance metadata to XCom, and optionally emits an Airflow Asset event. | ||
|
|
||
| The sensor can execute synchronously in standard worker slots or in deferrable mode using | ||
| :class:`~airflow.providers.databricks.triggers.databricks_delta_table.DatabricksDeltaTableVersionTrigger`. | ||
|
|
||
| .. note:: | ||
|
|
||
| **Scope (Phase 1 only):** This sensor provides task-level Delta table version observation inside a scheduled | ||
| DAG task run. A continuous external asset producer with long-running coalescing across runs is planned as a | ||
| follow-up Phase 2 feature. |
There was a problem hiding this comment.
This is use facing doc.
Users don't care about phases. We should just document what we currently have.
If explaining about future work has value then it's OK to mention it but then please explain why it's important to highlight? every feature of airflow can be enhanced further that is how software works :)
There was a problem hiding this comment.
I am not very happy with this doc. Feels very AI generated and not designed to be read by humans.
It feels as if the doc explains the code. parameter by parameter.
There are too many bold statements and too many things marked as important.
|
cc @moomindani for databricks team review |
|
Thanks @eladkal — I’ve updated the docs based on your feedback:
|
Add DatabricksDeltaTableVersionSensor and DatabricksDeltaTableVersionTrigger to monitor Delta table commit versions (Phase 1: scheduled-task sensor and deferrable trigger): - Rebase onto main and reuse UnityTableIdentity from airflow.providers.databricks.assets.databricks. - Validate Databricks connection workspace host against unity_table host case-insensitively. - Match Unity Catalog table parts case-insensitively, handling mixed-case table_name inputs. - Forward and serialize session_configuration, http_headers, client_parameters, hook_params, and query_tags into the deferrable trigger. - Fail fast on permanent query errors (missing table, permission denied) in deferrable trigger while retrying transient errors. - Restore baseline_version from trigger event in execute_complete before pushing provenance. - Emit precise task-success asset event extra metadata: observed Delta version and table identity (documenting that intermediate commits may be coalesced, and authors requiring pinned reads must VERSION AS OF). - Clarify in documentation that pinned reads require VERSION AS OF via ti.xcom_pull(task_ids="...", key="delta_table_version")["version"]. - Support baseline_version, target_version, and allow_recreation handling. - Support deferrable execution via async triggerer polling. - Add provider documentation including reproducibility notes and Phase 1 scope callout. - Register sensor and trigger in provider.yaml and get_provider_info.py.
- Rewrite delta_table_sensor.rst as standard how-to guide, removing Phase 1/2 notes and fixing typo - Register delta_table_sensor.rst under Databricks SQL how-to-guide in provider.yaml - Use compat conf and AirflowFailException / AirflowSensorTimeout in sensor - Narrow query result types in trigger and sensor to resolve MyPy errors - Defer private baseline version assignment until post-render execution - Use version-independent outlet_events accessor in test double
7152e22 to
a02cf8f
Compare
moomindani
left a comment
There was a problem hiding this comment.
Three findings, inline below. The first two I reproduced with airflow dags test (Airflow main, with only DatabricksSqlHook.run faked to return a version that goes up by one per call). The third I reproduced on a real workspace.
Separately, the docs should mention cost. Each poke runs a query on the SQL warehouse, so a poke_interval shorter than the warehouse's auto-stop keeps it running. This applies in deferrable mode too.
Drafted-by: Claude Code (Opus 5.5); reviewed by @moomindani before posting
|
|
||
| def _evaluate_version(self, current_version: int, table_name: str) -> bool: | ||
| """Evaluate if the current version meets sensor criteria.""" | ||
| if self._baseline_version is None and self.target_version is None: |
There was a problem hiding this comment.
With mode="reschedule" and no baseline_version, the sensor never succeeds. The enrolled baseline lives only on the operator instance. Every reschedule runs on a fresh instance, so each poke enrolls again and returns False. In my run the table committed 9 times (versions 5 to 13), the sensor logged Initialized baseline version to N on every poke, and it ended in AirflowSensorTimeout. A retry hits the same reset. Poke mode and deferrable mode are fine, because the baseline either stays in the process or is passed to the trigger.
Either keep the enrolled baseline across reschedules, or reject mode="reschedule" without baseline_version/target_version in __init__ so it fails loudly. A test that runs two separate sensor instances against increasing versions would cover it.
| return False | ||
|
|
||
| if self.target_version is not None: | ||
| if current_version >= self.target_version: |
There was a problem hiding this comment.
target_version is in template_fields, but unlike baseline_version it is never converted with int(). A templated value such as target_version="{{ params.v }}" renders to a string, and the task fails with TypeError: '>=' not supported between instances of 'int' and 'str'. The trigger receives the same string when deferred. Converting it to int in the same place as baseline_version (and testing it with a templated value) fixes both paths.
| return self.table_name | ||
|
|
||
| def _get_latest_version(self, table_name: str) -> tuple[int, str | None, str | None]: | ||
| sql = f"DESCRIBE HISTORY {table_name} LIMIT 1" |
There was a problem hiding this comment.
DESCRIBE HISTORY ... LIMIT 1 also returns maintenance commits, so the sensor fires when no data changed. On a real Unity Catalog managed table (Predictive Optimization ENABLE (inherited from METASTORE ...), which is the default for newer accounts), I ran two INSERTs and then OPTIMIZE and VACUUM. The history is 1 WRITE, 2 WRITE, 3 OPTIMIZE, 4 VACUUM START, 5 VACUUM END, so LIMIT 1 returns OPTIMIZE or VACUUM END as the "new version". With predictive optimization these commits arrive on their own schedule, so downstream Dags would run on compaction rather than on the external write that #74195 is about. The trigger (triggers/databricks_delta_table.py:131) has the same query.
Consider skipping maintenance operations (for example, reading a few history rows and taking the newest one whose operation is not OPTIMIZE/VACUUM START/VACUUM END), or at least making which operations count configurable and documenting the default.
closes: #74195
Description
Adds explicit Delta table version observation to the Databricks provider as Phase 1 of #74195.
Key additions:
**: Dataclass representing a static Unity Catalog table identity with validation and conversion to Airflow Asset (to_asset()`).: Sensor that queries Delta history (DESCRIBE HISTORY {table_name} LIMIT 1), detects new version commits, records provenance into XCom and Asset eventextra` metadata, and supports automatic asset outlet declaration.target_version, andallow_recreationhandling for tables dropped and recreated.providers/databricks/docs/operators/delta_table_sensor.rst.provider.yamlandget_provider_info.py.